Java NIO - 使用 netty 模拟 dubbo RPC 远程服务调用


使用 Netty 从零实现一个 RPC(类似 Dubbo)是一个非常经典的工程实践案例。为了遵循由简单到复杂的工程原则,我们将按以下五个步骤逐步进行深化。

  • Stage 1:跑通 Netty 最基础的通信,用反射完成一个最基本的调用。
  • Stage 2:引入 RequestId 和 CompletableFuture,完成单连接多路复用。
  • Stage 3:引入心跳和重连,保证了长连接的网络健壮性。
  • Stage 4:重构掉废弃 API,设计自定义二进制私有协议头和高性能 Kryo 序列化。
  • Stage 5:跳出单机思维,引入服务注册发现中心软负载均衡轮询算法。


基础案例

首先我们从最小可行性产品开始。在这个阶段,我们不考虑复杂的路由、服务注册中心(ZooKeeper/Nacos)、多线程模型优化和复杂的序列化,只聚焦于核心:如何通过 Netty 实现客户端发送一个请求对象,服务端接收、执行并返回结果。一个最基础的 RPC 包含以下四个要素:

  • API 接口:服务提供方和服务消费方共同约定的契约。
  • 网络传输(Netty):负责数据的发送与接收。
  • 协议与序列化:定义数据传输格式。为了简化直接使用 Netty 自带的 ObjectEncoder/Decoder。
  • 反射调用:服务端收到请求后,通过反射调用本地真正的方法。


接口和契约类

首先,定义一个服务接口,以及传输的请求/响应对象,传输对象必须实现 Serializable 接口。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
/**
* 定义服务接口
*/
public interface HelloService {
String sayHello(String msg);
}


/**
* 定义 RPC 请求体
*/
@Data
public class RpcRequest implements Serializable {
private String className;
private String methodName;
private Class<?>[] parameterTypes;
private Object[] parameters;
}


/**
* 定义 RPC 响应体
*/
@Data
public class RpcResponse implements Serializable {
private Object result;
private Throwable error;
}


服务端实现

服务端需要实现接口,并启动一个 Netty 服务监听端口。收到 RpcRequest 后,反射调用实现类,再把 RpcResponse 写回。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
/**
* 服务端实现接口
*/
public class HelloServiceImpl implements HelloService {
@Override
public String sayHello(String msg) {
return "【Server 回应】: 你好,我已经收到了你的消息 -> " + msg;
}
}


/**
* 服务端 Netty 业务处理器
*/
public class RpcServerHandler extends SimpleChannelInboundHandler<RpcRequest> {

@Override
protected void channelRead0(ChannelHandlerContext ctx, RpcRequest request) throws Exception {
RpcResponse response = new RpcResponse();
try {
// 极简模拟:直接硬编码创建实现类对象(后续会演进为从容器/注册表获取)
Object serviceBean = new HelloServiceImpl();

// 反射调用
Method method = serviceBean.getClass().getMethod(request.getMethodName(), request.getParameterTypes());
Object result = method.invoke(serviceBean, request.getParameters());

response.setResult(result);
} catch (Throwable t) {
response.setError(t);
}
// 写回给客户端
ctx.writeAndFlush(response);
}
}


/**
* 服务端 Netty 启动类
*/
public class RpcServer {
private int port;

public RpcServer(int port) {
this.port = port;
}

public void start() throws Exception {
EventLoopGroup bossGroup = new NioEventLoopGroup(1);
EventLoopGroup workerGroup = new NioEventLoopGroup();
try {
ServerBootstrap b = new ServerBootstrap();
b.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class)
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) throws Exception {
ChannelPipeline pipeline = ch.pipeline();
// 使用 Netty 自带的基于 JDK 序列化的编解码器(最简 case)
pipeline.addLast(new ObjectDecoder(Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null)));
pipeline.addLast(new ObjectEncoder());
// 自己的业务处理器
pipeline.addLast(new RpcServerHandler());
}
});
ChannelFuture f = b.bind(port).sync();
System.out.println("RPC Server 成功启动,监听端口: " + port);
f.channel().closeFuture().sync();
} finally {
bossGroup.shutdownGracefully();
workerGroup.shutdownGracefully();
}
}

/**
* 服务端启动
*/
public static void main(String[] args) throws Exception {
new RpcServer(8080).start();
}
}


客户端实现

客户端需要使用动态代理,让开发者像调用本地方法一样调用远程服务。同时利用 Netty 发送数据并同步阻塞等待结果。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
/**
* 客户端 Netty 处理器
*
*【在这个极简 RPC 客户端中,我们需要解决一个问题:异步转同步】
* 1. 当业务代码调用 helloService.sayHello("xxx") 时,它是同步的,必须拿到结果才能往下走。
* 2. 但 Netty 的网络通信是异步的。客户端把请求发出去后,业务线程就没事干了,而服务器的响应数据,是由 Netty 的 I/O 线程(RpcClientHandler 里的 channelRead0)接收到的。
* 3. 这里使用 SynchronousQueue,它是一个不存储元素的阻塞队列。业务线程发送完网络请求后,立刻调用 queue.take()。因为这时候结果还没回来,业务线程就在这里死等(阻塞)。
* 4. 过了几毫秒,Netty 收到服务器的响应,Netty 线程立刻调用 queue.put(response)。
* 5. 此时,两个线程在 SynchronousQueue 成功“接头”,Netty 线程把数据递过去,业务线程立刻被唤醒,拿着数据返回给上层代码。
*/
public class RpcClientHandler extends SimpleChannelInboundHandler<RpcResponse> {

// 使用阻塞队列来简单实现「异步转同步」获取结果
private final SynchronousQueue<RpcResponse> queue = new SynchronousQueue<>();

@Override
protected void channelRead0(ChannelHandlerContext ctx, RpcResponse msg) throws Exception {
queue.put(msg); // 收到服务器响应,注入队列。这一行会卡住,直到有另一个线程来执行 queue.take()。
}

public RpcResponse getResponse() throws InterruptedException {
return queue.take(); // 也会阻塞等待,直到拿到结果
}

@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
cause.printStackTrace();
ctx.close();
}
}


/**
* 客户端代理与连接管理器
*/
public class RpcClientProxy {

@SuppressWarnings("unchecked")
public static <T> T createProxy(Class<T> serviceClass, String host, int port) {
return (T) newProxyInstance(
serviceClass.getClassLoader(), // classloader
new Class<?>[]{serviceClass}, // interfaces
(proxy, method, args) -> { // handler

// 1. 组装请求体
RpcRequest request = new RpcRequest();
request.setClassName(serviceClass.getName());
request.setMethodName(method.getName());
request.setParameterTypes(method.getParameterTypes());
request.setParameters(args);

// 2. 启动 Netty 客户端发送数据(极简模式下每次调用新建连接,生产莫这样用)
EventLoopGroup group = new NioEventLoopGroup();
RpcClientHandler clientHandler = new RpcClientHandler();
try {
Bootstrap b = new Bootstrap();
b.group(group)
.channel(NioSocketChannel.class)
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) throws Exception {
ChannelPipeline pipeline = ch.pipeline();
pipeline.addLast(new ObjectEncoder());
pipeline.addLast(new ObjectDecoder(Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null)));
pipeline.addLast(clientHandler);
}
});

ChannelFuture f = b.connect(host, port).sync();
// 发送请求
f.channel().writeAndFlush(request).sync();

// 3. 阻塞等待服务器返回结果
RpcResponse response = clientHandler.getResponse();

if (response.getError() != null) {
throw response.getError();
}
return response.getResult();
} finally {
group.shutdownGracefully();
}
}
);
}

/**
* 服务调用端测试方法
*/
public static void main(String[] args) {
// 创建代理对象
HelloService helloService = RpcClientProxy.createProxy(HelloService.class, "127.0.0.1", 8080);

// 像调用本地方法一样调用远程服务
String result = helloService.sayHello("Hello Netty RPC!");
System.out.println("客户端收到远程调用结果: " + result);
}
}


测试验证

依次执行以下步骤:

  • 启动 RpcServer。
  • 运行 RpcClientProxy。

  • 你将看到控制台成功打印出通过网络反射调用后返回的字符串。


长连接复用与多服务路由

需要解决的问题

在 RPC 框架中,网络通信模型的演进(从短连接到长连接复用)才是最核心、最硬核的架构部分。这里我们需要实现长连接复用与多服务路由。它要解决的核心痛点如下:

  • 痛点一:短连接开销巨大。上面的基础案例,客户端每调用一次方法,就要 new NioEventLoopGroup(),建连、发数据、断连、销毁线程池,开销巨大。我们需要客户端与服务端建立一条(或几条)常驻的长连接,所有 RPC 请求都复用这个通道(Channel)发送。
  • 痛点二:高并发下的请求错乱(多路复用)。当多个线程复用同一个 Channel 发送请求时,Netty 是异步返回的。比如线程 A 发了请求 1,线程 B 发了请求 2;服务器可能先返回响应 2,再返回响应 1。客户端如何让线程 A 拿到响应 1,线程 B 拿到响应 2 就是个棘手的问题。这里我们需要引入请求 ID (RequestId) 和 全局期约(CompletableFuture) 映射表。
  • 痛点三:服务端无法支持多个服务。在基础案例中,Server 端写死了 new HelloServiceImpl()。这里我们在服务端引入一个本地服务注册表(Service Registry),用一个 Map 维护服务名与实现类的映射。


代码实现演进

为了实现多路复用,我们需要改造 RpcRequest 和 RpcResponse,加入一个 requestId。客户端发送请求时,把一个 CompletableFuture 存入全局 Map,然后业务线程阻塞在 future.get() 上。当 Netty 收到响应时,根据响应里的 requestId 找到对应的 Future 并赋值,从而唤醒业务线程。


契约升级

加入 RequestId:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
@Data
public class RpcRequest implements Serializable {
private String requestId; // 核心:唯一标识一次请求
private String className;
private String methodName;
private Class<?>[] parameterTypes;
private Object[] parameters;
}

@Data
public class RpcResponse implements Serializable {
private String requestId; // 核心:对应请求的 ID
private Object result;
private Throwable error;
}


服务端升级

服务端不再硬编码,而是提供一个注册方法,支持多服务路由。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
/**
* 服务端业务处理器(共享单例,处理所有路由)
*/
public class RpcServerHandler extends SimpleChannelInboundHandler<RpcRequest> {

// 关键:持有服务注册表
private final Map<String, Object> serviceMap;

public RpcServerHandler(Map<String, Object> serviceMap) {
this.serviceMap = serviceMap;
}

@Override
protected void channelRead0(ChannelHandlerContext ctx, RpcRequest request) throws Exception {
RpcResponse response = new RpcResponse();
// 传递请求 ID
response.setRequestId(request.getRequestId());
try {
// 根据接口名动态路由到实现类
Object serviceBean = serviceMap.get(request.getClassName());
if (serviceBean == null) {
throw new RuntimeException("未找到服务实现类: " + request.getClassName());
}
// 反射调用
Method method = serviceBean.getClass().getMethod(request.getMethodName(), request.getParameterTypes());
Object result = method.invoke(serviceBean, request.getParameters());
response.setResult(result);
} catch (Throwable t) {
response.setError(t);
}
ctx.writeAndFlush(response);
}
}

服务端启动类

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
/**
* 服务端 Netty 启动类
*/
public class RpcServer {
private int port;

public RpcServer(int port) {
this.port = port;
}

// 本地服务注册表
private final Map<String, Object> serviceMap = new ConcurrentHashMap<>();

// 发布服务的方法
public void registerService(Class<?> interfaceClass, Object serviceImpl) {
serviceMap.put(interfaceClass.getName(), serviceImpl);
}

public void start() throws Exception {
EventLoopGroup bossGroup = new NioEventLoopGroup(1);
EventLoopGroup workerGroup = new NioEventLoopGroup();
try {
ServerBootstrap b = new ServerBootstrap();
b.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class)
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) throws Exception {
ChannelPipeline pipeline = ch.pipeline();
// 使用 Netty 自带的基于 JDK 序列化的编解码器(最简 case)
pipeline.addLast(new ObjectDecoder(Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null)));
pipeline.addLast(new ObjectEncoder());
// 自己的业务处理器,传入注册表!
pipeline.addLast(new RpcServerHandler(serviceMap));
}
});
ChannelFuture f = b.bind(port).sync();
System.out.println("RPC Server 成功启动,监听端口: " + port);
f.channel().closeFuture().sync();
} finally {
bossGroup.shutdownGracefully();
workerGroup.shutdownGracefully();
}
}
}


客户端升级

客户端支持长连接与多路复用。客户端需要保持长连接,并使用 CompletableFuture 解决异步响应的匹配问题。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
/**
* 客户端 Netty 处理器
*/
public class RpcClientHandler extends SimpleChannelInboundHandler<RpcResponse> {

// 核心:维护 请求ID -> 期约Future 的映射表(全局共享)
private final Map<String, CompletableFuture<RpcResponse>> pendingResponseMap = new ConcurrentHashMap<>();

// 暴露给客户端代理,用来注册等待的 Future
public CompletableFuture<RpcResponse> registerFuture(String requestId) {
CompletableFuture<RpcResponse> future = new CompletableFuture<>();
pendingResponseMap.put(requestId, future);
return future;
}

@Override
protected void channelRead0(ChannelHandlerContext ctx, RpcResponse response) throws Exception {
// 收到响应,根据 requestId 找到对应的 Future
CompletableFuture<RpcResponse> future = pendingResponseMap.remove(response.getRequestId());
if (future != null) {
future.complete(response); // 完美的“精准唤醒”对应线程
}
}

@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
cause.printStackTrace();
ctx.close();
}
}

客户端连接管理器与代理(单例长连接)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
public class RpcClientContainer {
private final String host;
private final int port;
private Channel channel;
private RpcClientHandler clientHandler;
private final EventLoopGroup group = new NioEventLoopGroup();

public RpcClientContainer(String host, int port) {
this.host = host;
this.port = port;
}

// 初始化长连接(只在应用启动时连接一次)
public void init() throws Exception {
clientHandler = new RpcClientHandler();
Bootstrap b = new Bootstrap();
b.group(group)
.channel(NioSocketChannel.class)
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ChannelPipeline pipeline = ch.pipeline();
pipeline.addLast(new ObjectEncoder());
pipeline.addLast(new ObjectDecoder(Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null)));
pipeline.addLast(clientHandler);
}
});

ChannelFuture f = b.connect(host, port).sync();
this.channel = f.channel();
System.out.println("RPC 客户端长连接建立成功!");
}

// 关闭释放资源
public void close() {
group.shutdownGracefully();
}

// 创建代理对象
@SuppressWarnings("unchecked")
public <T> T createProxy(Class<T> serviceClass) {
return (T) Proxy.newProxyInstance(
serviceClass.getClassLoader(),
new Class<?>[]{serviceClass},
(proxy, method, args) -> {

// 1. 组装请求,生成唯一 ID
RpcRequest request = new RpcRequest();
request.setRequestId(UUID.randomUUID().toString());
request.setClassName(serviceClass.getName());
request.setMethodName(method.getName());
request.setParameterTypes(method.getParameterTypes());
request.setParameters(args);

// 2. 注册 Future
CompletableFuture<RpcResponse> future = clientHandler.registerFuture(request.getRequestId());

// 3. 复用同一个 channel 发送数据(不需要每次重连了!),这里首先做一下防御编程。
if (channel == null || !channel.isActive()) {
throw new RuntimeException("RPC 调用失败:当前网络连接已断开,正在尝试重连中...");
}
channel.writeAndFlush(request);

// 4. 同步阻塞等待这个 Future 的完成
RpcResponse response = future.get();

if (response.getError() != null) {
throw response.getError();
}
return response.getResult();
}
);
}
}


测试验证

增加一个新接口 OrderService,验证 “多服务路由”。

1
2
3
4
5
6
7
8
9
10
public interface OrderService {
String getOrderSymbol(String userId); // 根据用户ID获取最新订单号
}

public class OrderServiceImpl implements OrderService {
@Override
public String getOrderSymbol(String userId) {
return "【OrderServer】用户 " + userId + " 的最新订单号为:TX20260705001";
}
}

服务端和客户端的启动类:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
// ---- 服务端启动 ----
public class ServerMain {
public static void main(String[] args) throws Exception {
RpcServer server = new RpcServer(8080);

// 注册多个服务
server.registerService(HelloService.class, new HelloServiceImpl());
server.registerService(OrderService.class, new OrderServiceImpl());

server.start();
}
}

// ---- 客户端启动与多线程调用测试 ----
public class ClientMain {
public static void main(String[] args) throws Exception {
RpcClientContainer container = new RpcClientContainer("127.0.0.1", 8080);
container.init(); // 初始化长连接

// 获取不同的代理对象
HelloService helloService = container.createProxy(HelloService.class);
OrderService orderService = container.createProxy(OrderService.class);

// 模拟多线程、多服务并发调用,测试长连接多路复用
for (int i = 0; i < 3; i++) {
final int index = i;
new Thread(() -> {
String res1 = helloService.sayHello("线程-" + index);
System.out.println(res1);

String res2 = orderService.getOrderSymbol("User-" + index);
System.out.println(res2);
}).start();
}

// 维持主线程不退出
Thread.sleep(2000);
container.close();
}
}

客户端打印:

1
2
3
4
5
6
7
RPC 客户端长连接建立成功!
【Server 回应】: 你好,我已经收到了你的消息 -> 线程-0
【Server 回应】: 你好,我已经收到了你的消息 -> 线程-1
【Server 回应】: 你好,我已经收到了你的消息 -> 线程-2
【OrderServer】用户 User-0 的最新订单号为:TX20260705001
【OrderServer】用户 User-1 的最新订单号为:TX20260705001
【OrderServer】用户 User-2 的最新订单号为:TX20260705001

在这个阶段中,我们完成了 RPC 框架最核心的蜕变:用一条长连接 + Map,支撑起了并发多线程、多接口的远程调用。性能比基础案例提升了几个数量级。

不过,细心的你可能发现了另一个隐患:在 RpcClientContainer 中,我们所有的请求都共用一条 channel。如果高并发极端场景下,大量的读写数据导致这条 Channel 堵塞了(或者长连接因网络闪断意外关闭了),整个客户端就会陷入瘫痪。为了解决这个问题,在接下来的案例升级中,业界(包括 Dubbo)通常会引入连接池(Channel Pool) 或者心跳检测与断线重连机制来解决这个问题。


实现连接池与心跳重连机制

需要解决的问题

在 Stage 2 中,我们实现了一条长连接的多路复用。但如果在生产环境中,这条常驻的长连接因为网络闪断、运营商抖动或者长时间没有数据传输被防火墙切断,客户端的后续请求就会全部失败。为了让这个 RPC 框架达到 “工业级” 的健壮性,Stage 3 我们需要攻克两大核心机制:

  1. 心跳检测与自动重连(Keep-Alive & Reconnect):在空闲时发送微小的探测包,发现连接断开立刻自动重新连接。
  2. 连接池(Channel Pool):高并发下如果单条 Channel 成为瓶颈或意外受损,连接池可以提供多条通道分担压力,并提供自动剔除坏连接的能力。

Netty 提供了开箱即用的 IdleStateHandler(空闲状态处理器)。

  • 客户端:如果在一段时间内(比如 5 秒)没有向服务端发任何数据,就自动发送一个“心跳契约包”给服务端。
  • 服务端:如果长时间(比如 15 秒)没有收到客户端的任何请求或心跳,则认为客户端已掉线,主动关闭该 Channel 释放资源。
  • 重连:客户端如果检测到 Channel 断开了,触发监听器自动发起异步重连。


代码实现演进

契约升级

为了支持心跳,我们先定义一个轻量级的 “心跳乒乓包”,不走复杂的业务反射。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
@Data
public class RpcRequest implements Serializable {
private String requestId;
private String className;
private String methodName;
private Class<?>[] parameterTypes;
private Object[] parameters;

// True 表示是心跳探测包,无需反射调用业务
private boolean heartbeat;

// 行内生成一个标准心跳包
public static RpcRequest createHeartbeat() {
RpcRequest r = new RpcRequest();
r.setHeartbeat(true);
return r;
}
}


服务端升级

服务端开启空闲超时,踢掉死连接。使用 IdleStateHandler 监听读空闲。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
public class RpcServerHandler extends SimpleChannelInboundHandler<RpcRequest> {

// 关键:持有服务注册表
private final Map<String, Object> serviceMap;

public RpcServerHandler(Map<String, Object> serviceMap) {
this.serviceMap = serviceMap;
}

@Override
protected void channelRead0(ChannelHandlerContext ctx, RpcRequest request) throws Exception {
// 👉🏻 1. 如果是客户端发来的心跳包,直接无视或打印,不进行业务反射
if (request.isHeartbeat()) {
System.out.println("【Server】收到客户端心跳维持包...");
return;
}

// 2. 正常业务处理(同 Stage 2)
RpcResponse response = new RpcResponse();
// 传递请求 ID
response.setRequestId(request.getRequestId());
try {
// 根据接口名动态路由到实现类
Object serviceBean = serviceMap.get(request.getClassName());
if (serviceBean == null) {
throw new RuntimeException("未找到服务实现类: " + request.getClassName());
}
// 反射调用
Method method = serviceBean.getClass().getMethod(request.getMethodName(), request.getParameterTypes());
Object result = method.invoke(serviceBean, request.getParameters());
response.setResult(result);
} catch (Throwable t) {
response.setError(t);
}
ctx.writeAndFlush(response);
}

// 3. 核心:监听空闲事件
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception {
if (evt instanceof IdleStateEvent) {
System.out.println("【Server警告】长时间未收到客户端消息,判定为无效连接,断开。");
ctx.close(); // 直接关闭死连接
} else {
super.userEventTriggered(ctx, evt);
}
}
}

服务端 Pipeline 配置修改:

1
2
3
4
5
6
7
8
9
10
11
@Override
protected void initChannel(SocketChannel ch) throws Exception {
ChannelPipeline pipeline = ch.pipeline();

// 参数含义:15秒内没有读到客户端数据,触发读空闲事件
pipeline.addLast(new IdleStateHandler(15, 0, 0, TimeUnit.SECONDS));

pipeline.addLast(new ObjectDecoder(Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null)));
pipeline.addLast(new ObjectEncoder());
pipeline.addLast(new RpcServerHandler(serviceMap));
}


客户端升级

客户端需要在空闲时发心跳,并且在连接断开时,利用 Netty 的 ChannelFutureListener 发起自动重连。

  • 初次启动守住大门:当你在 ClientMain 中调用 container.init() 时,主线程会卡在 future.sync() 这一行,直到网络通道彻底打通、this.channel 被成功赋值,主线程才会被放行去创建代理和启动多线程测试。
  • 重连不影响线程池:当后面网络意外断开时,触发 channelInactive,此时我们调用 doConnect(false) 走异步分支。在重连成功的几秒钟空窗期内,如果业务线程来调用,会被我们的 “健壮性补丁” 拦截,抛出可控的业务异常,而不会报底层无厘头的 NullPointerException。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
public class RpcClientContainer {
private final String host;
private final int port;
private Channel channel;
private RpcClientHandler clientHandler;
private Bootstrap bootstrap;
private final EventLoopGroup group = new NioEventLoopGroup();

public RpcClientContainer(String host, int port) {
this.host = host;
this.port = port;
}

// 初始化长连接(只在应用启动时连接一次)
public void init() throws Exception {
clientHandler = new RpcClientHandler();
bootstrap = new Bootstrap();
bootstrap.group(group)
.channel(NioSocketChannel.class)
.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 5000) // 👉🏻 5秒建连超时
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) {
ChannelPipeline pipeline = ch.pipeline();

// 👉🏻 5秒内没有向服务器发送过数据(写空闲),就触发事件(readerIdleTime、writeIdleTime、allIdleTime、时间单位)
pipeline.addLast(new IdleStateHandler(0, 5, 0, TimeUnit.SECONDS));

pipeline.addLast(new ObjectEncoder());
pipeline.addLast(new ObjectDecoder(Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null)));

// 👉🏻 处理业务及心跳触发的 Handler
pipeline.addLast(new RpcClientHandlerInner());
}
});

// 初次连接:使用同步等待,确保 channel 初始化完成才结束 init()
doConnect(true);
}

// 👉🏻 封装核心建连逻辑,支持同步初次连接和异步后续重连
private void doConnect(boolean isInit) {
if (channel != null && channel.isActive()) {
return;
}
System.out.println("尝试连接 RPC 服务器 [" + host + ":" + port + "]...");
try {
ChannelFuture future = bootstrap.connect(host, port);
if (isInit) {
// 如果是初次初始化,强行同步阻塞死等 5 秒建连,成功后才放行
future.sync();
this.channel = future.channel();
System.out.println("RPC 客户端初次连接成功!");
} else {
// 如果是后续的断线自动重连,依然走异步监听监听器,不阻塞 Netty 的 EventLoop 线程
future.addListener((ChannelFutureListener) f -> {
if (f.isSuccess()) {
this.channel = f.channel();
System.out.println("RPC 客户端自动重连成功!");
} else {
System.out.println("RPC 客户端重连失败,3秒后准备再次重连...");
f.channel().eventLoop().schedule(() -> doConnect(false), 3, TimeUnit.SECONDS);
}
});
}
} catch (Exception e) {
if (isInit) {
System.err.println("RPC 客户端初次连接服务器失败!原因: " + e.getMessage());
// 初次连接都失败了,直接抛出异常阻止程序盲目启动
throw new RuntimeException("RPC 客户端初始化失败,无法连接服务器", e);
}
}
}

// 关闭释放资源
public void close() {
group.shutdownGracefully();
}

// 👉🏻 内部 Handler:负责拦截空闲事件并发送心跳,以及处理连接断开时的重连
private class RpcClientHandlerInner extends SimpleChannelInboundHandler<RpcResponse> {
@Override
protected void channelRead0(ChannelHandlerContext ctx, RpcResponse msg) throws Exception {
// 桥接到全局的响应中心
clientHandler.channelRead0(ctx, msg);
}

// 👉🏻 写空闲触发时,发送心跳包
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception {
if (evt instanceof IdleStateEvent) {
// 发送心跳维持包
ctx.writeAndFlush(RpcRequest.createHeartbeat());
} else {
super.userEventTriggered(ctx, evt);
}
}

// 核心:连接一旦意外断开,立刻触发重连
@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception {
System.out.println("【客户端警告】检测到连接断开,触发自动重连机制!");
super.channelInactive(ctx);
// 借用 Netty 的线程池异步执行重连逻辑。传入 false,绝不能在 EventLoop 里执行 sync() 导致死锁
ctx.channel().eventLoop().schedule(() -> doConnect(false), 3, TimeUnit.SECONDS);
}
}


// 创建代理对象(保持不变)
@SuppressWarnings("unchecked")
public <T> T createProxy(Class<T> serviceClass) {
return (T) Proxy.newProxyInstance(
serviceClass.getClassLoader(),
new Class<?>[]{serviceClass},
(proxy, method, args) -> {

// 1. 组装请求,生成唯一 ID
RpcRequest request = new RpcRequest();
request.setRequestId(UUID.randomUUID().toString());
request.setClassName(serviceClass.getName());
request.setMethodName(method.getName());
request.setParameterTypes(method.getParameterTypes());
request.setParameters(args);

// 2. 注册 Future
CompletableFuture<RpcResponse> future = clientHandler.registerFuture(request.getRequestId());

// 3. 复用同一个 channel 发送数据(不需要每次重连了!)
if (channel == null || !channel.isActive()) {
throw new RuntimeException("RPC 调用失败:当前网络连接已断开,正在尝试重连中...");
}
channel.writeAndFlush(request);

// 4. 同步阻塞等待这个 Future 的完成
RpcResponse response = future.get();

if (response.getError() != null) {
throw response.getError();
}
return response.getResult();
}
);
}
}


测试健壮性

微调 ClientMain:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
public class ClientMain {

public static void main(String[] args) throws Exception {
RpcClientContainer container = new RpcClientContainer("127.0.0.1", 8080);
container.init(); // 初始化长连接

// 获取不同的代理对象
HelloService helloService = container.createProxy(HelloService.class);
OrderService orderService = container.createProxy(OrderService.class);

// 模拟多线程、多服务并发调用,测试长连接多路复用
for (int i = 0; i < 3; i++) {
final int index = i;
new Thread(() -> {
String res1 = helloService.sayHello("线程-" + index);
System.out.println(res1);

String res2 = orderService.getOrderSymbol("User-" + index);
System.out.println(res2);
}).start();
}

// 维持主线程不退出
/*Thread.sleep(2000);
container.close();*/
}
}

依次进行执行验证:

  • 正常启动 ServerMain 和 ClientMain。
  • 客户端开始发送请求。此时你可以看到服务端定期打印:【Server】收到客户端心跳维持包…
  • 断网模拟:直接在控制台强行把 ServerMain 关闭(Kill 掉)。
  • 客户端控制台会立刻报警:【客户端警告】检测到连接断开,触发自动重连机制! 并每隔 3 秒持续打印 尝试连接 RPC 服务器…
  • 恢复模拟:重新把 ServerMain 启动起来。
  • 客户端会在几秒内捕捉到服务的复活,打印 RPC 客户端连接成功!,整条链路完好如初,不需要重启客户端!


自定义协议与高性能序列化升级

自定义私有协议

在这一阶段,我们将彻底告别已经被废弃且漏洞百出的 ObjectDecoder/Encoder。我们将像 Dubbo 一样,设计一套专属的自定义二进制私有协议。我们设计的协议头共占 8 个字节,结构如下:

1
2
3
4
5
6
+---------------------------------------------------------------+
| 魔数 (4B) | 序列化类型 (1B) | 消息类型 (1B) | 数据长度 (2B) |
+---------------------------------------------------------------+
| 真正的数据 (Payload) |
+---------------------------------------------------------------+

  • 魔数 (Magic Number) - 4字节:固定为 0xCAFEBABE(或者你喜欢的任意 4 字节)。用来筛选和过滤非本框架的非法请求,防止恶意攻击。
  • 序列化类型 (Serialization Type) - 1字节:标记使用的是哪种序列化方式。例如:0x01 代表 Kryo,0x02 代表 JSON。这为以后的多序列化扩展打下基础。
  • 消息类型 (MessageType) - 1字节:标记这是一个请求、响应还是心跳包。例如:0x01 请求,0x02 响应,0x03 心跳。这样心跳包就不需要复杂的业务反序列化了。
  • 数据长度 (Length) - 2字节:标记后面真正的数据体(Payload)占用了多少字节。最大支持 $2^{16}-1 = 65535$ 字节(如果需要传大对象,可改为 4 字节)。


引入高性能序列化组件 Kryo

Kryo 是一个快速、高效的 Java 对象图形序列化框架,性能和压缩率远超 JDK 原生序列化。

1
2
3
4
5
<dependency>
<groupId>com.esotericsoftware</groupId>
<artifactId>kryo</artifactId>
<version>5.5.0</version>
</dependency>


Kryo 工具类封装

由于 Kryo 不是线程安全的,必须使用 ThreadLocal 来保证每个线程拥有独立的 Kryo 实例。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
import com.esotericsoftware.kryo.Kryo;
import com.esotericsoftware.kryo.io.Input;
import com.esotericsoftware.kryo.io.Output;
import com.owlias.janus.comet.client.zdemo.common.RpcRequest;
import com.owlias.janus.comet.client.zdemo.common.RpcResponse;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;

public class KryoSerializer {
private static final ThreadLocal<Kryo> kryoThreadLocal = ThreadLocal.withInitial(() -> {
Kryo kryo = new Kryo();
// 显式注册我们需要传输 include 的类,提高安全性和效率
kryo.register(RpcRequest.class);
kryo.register(RpcResponse.class);
kryo.register(Class[].class);
kryo.register(Class.class);
kryo.register(Object[].class);
kryo.setReferences(true); // 支持循环引用
kryo.setRegistrationRequired(false); // 关闭强行注册要求
return kryo;
});

public static byte[] serialize(Object obj) {
Kryo kryo = kryoThreadLocal.get();
ByteArrayOutputStream out = new ByteArrayOutputStream();
Output output = new Output(out);
kryo.writeClassAndObject(output, obj);
output.close();
return out.toByteArray();
}

public static Object deserialize(byte[] bytes) {
Kryo kryo = kryoThreadLocal.get();
ByteArrayInputStream in = new ByteArrayInputStream(bytes);
Input input = new Input(in);
Object obj = kryo.readClassAndObject(input);
input.close();
return obj;
}
}


编写自定义编解码器

编码器:将对象打包成二进制流(Outbound)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
import com.owlias.janus.comet.client.zdemo.common.RpcRequest;
import com.owlias.janus.comet.client.zdemo.serialize.KryoSerializer;
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.MessageToByteEncoder;

/**
* 自定义编码器:将对象打包成二进制流(Outbound)
*/
public class MyProtocolEncoder extends MessageToByteEncoder<Object> {
private static final int MAGIC_NUMBER = 0xCAFEBABE;

@Override
protected void encode(ChannelHandlerContext ctx, Object msg, ByteBuf out) throws Exception {
// 1. 写入魔数 (4B)
out.writeInt(MAGIC_NUMBER);
// 2. 写入序列化类型 (1B) -> 0x01 代表 Kryo
out.writeByte((byte) 0x01);

// 3. 写入消息类型 (1B)
if (msg instanceof RpcRequest) {
RpcRequest req = (RpcRequest) msg;
if (req.isHeartbeat()) {
out.writeByte((byte) 0x03); // 心跳包
} else {
out.writeByte((byte) 0x01); // 业务请求
}
} else {
out.writeByte((byte) 0x02); // 业务响应
}

// 4. 序列化 Payload
byte[] bodyBytes = KryoSerializer.serialize(msg);

// 5. 写入数据长度 (2B) 并写入真正的数据
out.writeShort(bodyBytes.length);
out.writeBytes(bodyBytes);
}
}

解码器:解决粘包半包,剥离协议头(Inbound)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
import com.owlias.janus.comet.client.zdemo.common.RpcRequest;
import com.owlias.janus.comet.client.zdemo.serialize.KryoSerializer;
import io.netty.buffer.ByteBuf;
import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.LengthFieldBasedFrameDecoder;

/**
* 自定义解码器:解决粘包半包,剥离协议头(Inbound)
* 我们继承 LengthFieldBasedFrameDecoder(长度域解码器),让 Netty 自动帮我们攒够一个完整的数据包再交给我们处理。
*/
public class MyProtocolDecoder extends LengthFieldBasedFrameDecoder {
private static final int MAGIC_NUMBER = 0xCAFEBABE;
private static final int HEADER_LENGTH = 8; // 4 + 1 + 1 + 2

public MyProtocolDecoder() {
// 参数含义:最大帧长65535,长度域偏移量为6(前4B魔数+1B序列化+1B类型),长度域占用2B,长度调整为0,跳过0字节
super(65535, 6, 2, 0, 0);
}

@Override
protected Object decode(ChannelHandlerContext ctx, ByteBuf in) throws Exception {
// 使用 Netty 自带的长度拆包器进行第一层拆包
ByteBuf frame = (ByteBuf) super.decode(ctx, in);
if (frame == null) return null;

try {
// 1. 校验魔数
int magic = frame.readInt();
if (magic != MAGIC_NUMBER) {
throw new IllegalArgumentException("非法协议魔数: " + Integer.toHexString(magic));
}
// 2. 读取序列化类型
byte serializeType = frame.readByte();
// 3. 读取消息类型
byte messageType = frame.readByte();
// 4. 读取数据长度
int length = frame.readShort();

// 如果是快速心跳包,直接拦截组装并返回,无需反序列化
if (messageType == (byte) 0x03) {
return RpcRequest.createHeartbeat();
}

// 5. 解包出真正的 Payload 字节数组
byte[] bodyBytes = new byte[length];
frame.readBytes(bodyBytes);

// 6. Kryo 反序列化回对象
return KryoSerializer.deserialize(bodyBytes);
} finally {
frame.release(); // 记得释放 ByteBuf 内存
}
}
}


替换原编解码器

编解码器写好后,我们在服务端和客户端的 initChannel 中,把之前废弃的 ObjectDecoder/Encoder 移出队列,换上我们帅气的专属私有协议。

服务端 RpcServer 修改:

1
2
3
4
pipeline.addLast(new IdleStateHandler(15, 0, 0, TimeUnit.SECONDS));
pipeline.addLast(new MyProtocolDecoder()); // 自定义二进制解码
pipeline.addLast(new MyProtocolEncoder()); // 自定义二进制编码
pipeline.addLast(new RpcServerHandler(serviceMap));

客户端 RpcClientContainer 修改:

1
2
3
4
pipeline.addLast(new IdleStateHandler(0, 5, 0, TimeUnit.SECONDS));
pipeline.addLast(new MyProtocolEncoder()); // 自定义二进制编码
pipeline.addLast(new MyProtocolDecoder()); // 自定义二进制解码
pipeline.addLast(new RpcClientHandlerInner());


测试验证

按照阶段3测试,我们依然能够正常稳定运行。下面是客户端部分日志:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
尝试连接 RPC 服务器 [127.0.0.1:8080]...
RPC 客户端初次连接成功!
【Server 回应】: 你好,我已经收到了你的消息 -> 线程-0
【Server 回应】: 你好,我已经收到了你的消息 -> 线程-1
【Server 回应】: 你好,我已经收到了你的消息 -> 线程-2
【OrderServer】用户 User-0 的最新订单号为:TX20260705001
【OrderServer】用户 User-1 的最新订单号为:TX20260705001
【OrderServer】用户 User-2 的最新订单号为:TX20260705001
【客户端警告】检测到连接断开,触发自动重连机制!
尝试连接 RPC 服务器 [127.0.0.1:8080]...
RPC 客户端重连失败,3秒后准备再次重连...
尝试连接 RPC 服务器 [127.0.0.1:8080]...
RPC 客户端重连失败,3秒后准备再次重连...
尝试连接 RPC 服务器 [127.0.0.1:8080]...
RPC 客户端自动重连成功!


服务发现与软负载均衡

需要解决的问题

在前面的阶段中,我们的客户端(Consumer)在 RpcClientContainer 里把服务端的 IP 和端口(127.0.0.1:8080)给硬编码死了。但在实际的生产分布式系统中,服务端往往是一个集群。当某个服务端节点宕机或者新加了机器时,客户端必须能够动态感知,并将请求均匀地分摊到各个健康的节点上。这就是服务注册中心与负载均衡的核心价值。

为了遵循由浅入深的原则,我们不急于立刻引入复杂的 ZooKeeper 或 Nacos 外部集群,而是先在本地使用一个标准的接口设计,抽象出一个 虚拟的/轻量级的分布式注册中心模型。

  • 服务提供者 (Provider):启动时,将自己提供的接口名和自己的 IP:Port 注册到注册中心。
  • 服务消费者 (Consumer):调用方法前,先去注册中心“拉取”该接口对应的全部可用地址列表(如 [192.168.1.10:8080, 192.168.1.11:8080])。
  • 负载均衡 (Load Balance):客户端从地址列表中,通过一定的算法(如轮询 RoundRobin 或随机 Random)挑选出一个地址,然后发起网络调用。


代码实现演进

注册与负载接口

首先,我们需要解耦,定义好服务注册、发现以及负载均衡的通用标准。

1
2
3
4
5
6
7
8
9
10
11
12
13
// 1. 注册中心通用契约
public interface ServiceRegistry {
// 服务端调用:注册服务地址
void register(String serviceName, String serverAddress);
// 客户端调用:根据服务名发现所有的可用地址
List<String> discover(String serviceName);
}

// 2. 负载均衡算法契约
public interface LoadBalance {
// 从一堆服务地址中,挑选出一个
String select(List<String> addressList);
}


注册中心实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
/**
* 实现轻量级注册中心(本地模拟分布式感知)
* 这里我们用一个全局单例的 ConcurrentHashMap 模拟分布式注册中心。
* 真实场景下只需将这里替换为 ZooKeeper 的临时节点或 Nacos 的 Instance 注册即可,API 结构完全一致。
*/
public class LocalMockRegistry implements ServiceRegistry {
// 模拟分布式中心的数据存储 Map<服务名, 地址列表>
private static final Map<String, List<String>> registryCenter = new ConcurrentHashMap<>();

@Override
public void register(String serviceName, String serverAddress) {
registryCenter.computeIfAbsent(serviceName, k -> new CopyOnWriteArrayList<>()).add(serverAddress);
System.out.println("【注册中心】服务 [" + serviceName + "] 成功注册新节点 -> " + serverAddress);
}

@Override
public List<String> discover(String serviceName) {
return registryCenter.getOrDefault(serviceName, new ArrayList<>());
}
}


负载均衡器实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
/**
* 负载均衡器实现,以最常用的轮询算法为例:
* 轮询算法(Round Robin)需要保证在高并发下,轮流分配地址。我们使用 AtomicInteger 来保证线程安全。
*/
public class RoundRobinLoadBalance implements LoadBalance {
private final AtomicInteger index = new AtomicInteger(0);

@Override
public String select(List<String> addressList) {
if (addressList == null || addressList.isEmpty()) {
return null;
}
// 防止 Integer 越界变成负数,进行绝对值取模
int pos = Math.abs(index.getAndIncrement() % addressList.size());
return addressList.get(pos);
}
}


客户端升级

现在,客户端 RpcClientContainer 不再只盯着一个 host:port 了,它需要在运行时,根据你要调用的接口,动态去注册中心找地址!为了不频繁建连,我们内部通过一个 Map<String, Channel> 维护针对不同服务器地址的长连接。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
public class RpcClientRouter {
private final ServiceRegistry registry;
private final LoadBalance loadBalance = new RoundRobinLoadBalance();

// 核心:地址 -> 对应的长连接 Channel(连接池映射)
private final Map<String, Channel> channelMap = new ConcurrentHashMap<>();
private final EventLoopGroup group = new NioEventLoopGroup();
private final RpcClientHandler clientHandler = new RpcClientHandler(); // 共享的响应中心

public RpcClientRouter(ServiceRegistry registry) {
this.registry = registry;
}

// 根据动态获取的地址建立或复用长连接
private Channel getChannel(String address) throws Exception {
if (channelMap.containsKey(address)) {
Channel ch = channelMap.get(address);
if (ch.isActive()) return ch;
}

String[] parts = address.split(":");
String host = parts[0];
int port = Integer.parseInt(parts[1]);

Bootstrap b = new Bootstrap();
b.group(group).channel(NioSocketChannel.class)
.handler(new ChannelInitializer<>() {
@Override
protected void initChannel(Channel ch) {
// 依然沿用 Stage 4 的自定义私有协议编解码
ch.pipeline().addLast(new MyProtocolEncoder());
ch.pipeline().addLast(new MyProtocolDecoder());
ch.pipeline().addLast(clientHandler);
}
});

ChannelFuture f = b.connect(host, port).sync();
Channel channel = f.channel();
channelMap.put(address, channel);
return channel;
}

@SuppressWarnings("unchecked")
public <T> T createProxy(Class<T> serviceClass) {
return (T) Proxy.newProxyInstance(
serviceClass.getClassLoader(),
new Class<?>[]{serviceClass},
(proxy, method, args) -> {
String serviceName = serviceClass.getName();

// 1. 动态发现:去注册中心抓取这个服务的全部健康节点
List<String> addresses = registry.discover(serviceName);
if (addresses.isEmpty()) {
throw new RuntimeException("RPC 错误:没有找到可用的服务节点 [" + serviceName + "]");
}

// 2. 负载均衡:利用轮询算法挑出一个节点的地址
String selectedAddress = loadBalance.select(addresses);
System.out.println("负载均衡 selectedAddress:" + selectedAddress);

// 3. 路由连接:获取对应地址的长连接
Channel channel = getChannel(selectedAddress);

// 4. 发送请求(组装、Future等待等逻辑同 Stage 2/3)
RpcRequest request = new RpcRequest();
request.setRequestId(UUID.randomUUID().toString());
request.setClassName(serviceName);
request.setMethodName(method.getName());
request.setParameterTypes(method.getParameterTypes());
request.setParameters(args);

var future = clientHandler.registerFuture(request.getRequestId());
channel.writeAndFlush(request);

return future.get().getResult();
}
);
}
}

由于这里的 clientHandler 是一个共享的对象 ,它允许同一个实例被同时配置到多个 Channel 中,这种情况下你必须声明它是可被共享的,否则当你两次 rpcClientRouter.getChannel 就会报错:

1
com.xxx.client.zdemo.client.RpcClientHandler is not a @Sharable handler, so can't be added or removed multiple times.

声明 RpcClientHandler 是共享的(它实际上是无状态的):

1
2
3
4
@ChannelHandler.Sharable  ////
public class RpcClientHandler extends SimpleChannelInboundHandler<RpcResponse> {
// ...
}


测试验证

为了方便测试,我们做如下新增或修改。

第一,RpcServer 增加 异步启动方法:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
/**
* 异步启动方法
*/
public void startAsync() {
// 新开一个单独的线程,专门用来初始化和启动 Netty
Thread serverThread = new Thread(() -> {
try {
// 直接复用并调用原本的同步启动逻辑
this.start();
} catch (Exception e) {
System.err.println("【Server 错误】端口 " + port + " 异步启动失败: " + e.getMessage());
}
}, "RpcServer-Binder-Thread-" + port); // 给线程起个清爽的名字
serverThread.start(); // 启动线程,立即返回,不阻塞主线程
}

第二,新增服务端和客户端测试类:

ClusterServerMain

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
public class ClusterServerMain {
public static void main(String[] args) throws Exception {
ServiceRegistry registry = new LocalMockRegistry();

// 1. 模拟启动 8080 端口的服务器,并发布服务
RpcServer server1 = new RpcServer(8080);
server1.registerService(HelloService.class, new HelloServiceImpl());
server1.startAsync(); // 异步启动避免阻塞
registry.register(HelloService.class.getName(), "127.0.0.1:8080");

// 2. 模拟启动 8081 端口的服务器,集群部署
RpcServer server2 = new RpcServer(8081);
server2.registerService(HelloService.class, new HelloServiceImpl());
server2.startAsync();
registry.register(HelloService.class.getName(), "127.0.0.1:8081");
}
}

ClusterClientMain

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
public class ClusterClientMain {
public static void main(String[] args) throws Exception {
ServiceRegistry registry = new LocalMockRegistry();
registry.register(HelloService.class.getName(), "127.0.0.1:8080");
registry.register(HelloService.class.getName(), "127.0.0.1:8081");
RpcClientRouter router = new RpcClientRouter(registry);

HelloService helloService = router.createProxy(HelloService.class);

// 连续调用 4 次,观察两台服务器是否交替收到请求(轮询负载均衡)
for (int i = 1; i <= 4; i++) {
String result = helloService.sayHello("集群测试 " + i);
System.out.println("客户端收到结果: " + result);
}
}
}

分别启动服务端和客户端:

服务端日志:

1
2
3
4
【注册中心】服务 [com.xxx.client.zdemo.common.HelloService] 成功注册新节点 -> 127.0.0.1:8080
【注册中心】服务 [com.xxx.client.zdemo.common.HelloService] 成功注册新节点 -> 127.0.0.1:8081
RPC Server 成功启动,监听端口: 8080
RPC Server 成功启动,监听端口: 8081

客户端日志:

1
2
3
4
5
6
7
8
9
10
【注册中心】服务 [com.xxx.client.zdemo.common.HelloService] 成功注册新节点 -> 127.0.0.1:8080
【注册中心】服务 [com.xxx.client.zdemo.common.HelloService] 成功注册新节点 -> 127.0.0.1:8081
负载均衡 selectedAddress:127.0.0.1:8080
客户端收到结果: 【Server 回应】: 你好,我已经收到了你的消息 -> 集群测试 1
负载均衡 selectedAddress:127.0.0.1:8081
客户端收到结果: 【Server 回应】: 你好,我已经收到了你的消息 -> 集群测试 2
负载均衡 selectedAddress:127.0.0.1:8080
客户端收到结果: 【Server 回应】: 你好,我已经收到了你的消息 -> 集群测试 3
负载均衡 selectedAddress:127.0.0.1:8081
客户端收到结果: 【Server 回应】: 你好,我已经收到了你的消息 -> 集群测试 4